ποΈGitΠ―ΡΠ°ποΈ
Node / meshtastic / Meshtastic-Android / files / core / takserver / src / jvmAndroidMain / kotlin / org / meshtastic / core / takserver / TAKClientConnection.kt
Displaying Raw β’ Download
core/takserver/src/jvmAndroidMain/kotlin/org/meshtastic/core/takserver/TAKClientConnection.kt a5e6894fe8934aefce5036458407a69fc96ca541 (a5e6894f) Text, 15.27 KB
T8b949e/*
* Copyright (c) 2026 Meshtastic LLC
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with this program. If not, see <https://www.gnu.org/licenses/>.
*/
Tf0883e@fileTb4b4b4:Te6edf3SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooManyFunctionsTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72package T7ee787org.meshtastic.core.takserver
Tff7b72import T7ee787co.touchlab.kermit.Logger
Tff7b72import T7ee787kotlinx.coroutines.CancellationException
Tff7b72import T7ee787kotlinx.coroutines.CoroutineDispatcher
Tff7b72import T7ee787kotlinx.coroutines.CoroutineScope
Tff7b72import T7ee787kotlinx.coroutines.Job
Tff7b72import T7ee787kotlinx.coroutines.SupervisorJob
Tff7b72import T7ee787kotlinx.coroutines.cancel
Tff7b72import T7ee787kotlinx.coroutines.launch
Tff7b72import T7ee787kotlinx.coroutines.sync.Mutex
Tff7b72import T7ee787kotlinx.coroutines.sync.withLock
Tff7b72import T7ee787kotlinx.coroutines.withContext
Tff7b72import T7ee787java.io.BufferedOutputStream
Tff7b72import T7ee787java.io.InputStream
Tff7b72import T7ee787java.io.OutputStream
Tff7b72import T7ee787java.net.Socket
Tff7b72import T7ee787java.util.concurrent.atomic.AtomicBoolean
Tff7b72import T7ee787kotlin.concurrent.Volatile
Tff7b72import T7ee787kotlin.random.Random
Tff7b72import T7ee787kotlin.time.Clock
Tff7b72import T7ee787kotlin.time.Duration.Companion.milliseconds
Tff7b72import T7ee787kotlin.time.Instant
Tff7b72import T7ee787kotlinx.coroutines.isActive Tff7b72as Te6edf3coroutineIsActive
T8b949e/**
* Per-client state machine for a connected TAK client (ATAK / iTAK / WinTAK).
*
* This is the jvmAndroidMain implementation, using plain `java.net.Socket` (which is also the base class of
* [javax.net.ssl.SSLSocket] from [TAKServerJvm]) with blocking `InputStream`/`OutputStream` I/O wrapped in
* [Dispatchers.IO] coroutines.
*
* Responsibilities:
* - TAK protocol negotiation handshake (`t-x-takp-v` / `-q` / `-r`)
* - Read loop that frames `<event>` elements off the stream via [CoTXmlFrameBuffer]
* - Keepalive loop that emits a `t-x-d-d` event every [TAK_KEEPALIVE_INTERVAL_MS]
* - Serializing writes under a mutex so interleaved broadcasts never corrupt the XML stream
* - Lifecycle reporting up to [TAKServerJvm] via [onEvent] (`Connected`, `Disconnected`, `Error`, `ClientInfoUpdated`,
* `Message`)
*/
Tff7b72internal Tff7b72class T56d364TAKClientConnectionTb4b4b4(
Tff7b72private Tff7b72val Te6edf3socketTb4b4b4: Te6edf3SocketTb4b4b4,
Tff7b72val Te6edf3clientInfoTb4b4b4: Te6edf3TAKClientInfoTb4b4b4,
Tff7b72private Tff7b72val Te6edf3onEventTb4b4b4: Tb4b4b4(Te6edf3TAKConnectionEventTb4b4b4) Tff7b72-Tff7b72> Tffa657UnitTb4b4b4,
Tff7b72private Tff7b72val Te6edf3scopeTb4b4b4: Te6edf3CoroutineScopeTb4b4b4,
Tff7b72private Tff7b72val Te6edf3ioDispatcherTb4b4b4: Te6edf3CoroutineDispatcherTb4b4b4,
Tb4b4b4) Tb4b4b4{
Tff7b72private Tff7b72var Te6edf3currentClientInfo Tff7b72= Te6edf3clientInfo
Tff7b72private Tff7b72val Te6edf3frameBuffer Tff7b72= Te6edf3CoTXmlFrameBufferTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3inputStreamTb4b4b4: Te6edf3InputStream Tff7b72= Te6edf3socketTb4b4b4.Te6edf3getInputStreamTb4b4b4(Tb4b4b4)
T8b949e// Wrap the OutputStream in a BufferedOutputStream so that multiple small writes
T8b949e// (we emit a full XML event per write) coalesce into one syscall; flush() after
T8b949e// each event to push the bytes through TLS immediately.
Tff7b72private Tff7b72val Te6edf3outputStreamTb4b4b4: Te6edf3OutputStream Tff7b72= Te6edf3BufferedOutputStreamTb4b4b4(Te6edf3socketTb4b4b4.Te6edf3getOutputStreamTb4b4b4(Tb4b4b4)Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3writeMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
T8b949e/**
* Per-connection child scope. Every coroutine this class launches β the read loop, the keepalive loop, and every
* single send β is attached to [connectionScope] so that [emitDisconnected] can tear the whole connection down with
* one `connectionScope.cancel()`.
*
* Why this is critical: [broadcast] in [TAKServerJvm] fires `connection.send()` on **every** connected client for
* **every** CoT event coming off the mesh (and with a 56-node nodeDB each `nodeDBbyNum` emission fans out to ~56
* broadcasts). If [sendXml] launched those writes on the server-level [scope] β as the previous implementation did
* β a single dead connection could accumulate hundreds of in-flight write coroutines before it was removed from
* [TAKServerJvm.connections], and every one of them would spin up, hit the closed TLS socket, and log
* `SocketException: Socket closed` from `BufferedOutputStream.flush()`. Scoping writes to [connectionScope] means
* cancelling the scope wipes the entire backlog.
*
* Uses a [SupervisorJob] child of [scope]'s job so a single write failure doesn't cascade-cancel other connections
* on the same server.
*/
Tff7b72private Tff7b72val Te6edf3connectionScopeTb4b4b4: Te6edf3CoroutineScope Tff7b72=
Te6edf3CoroutineScopeTb4b4b4(Te6edf3SupervisorJobTb4b4b4(Te6edf3scopeTb4b4b4.Te6edf3coroutineContextTff7b72[Te6edf3JobTff7b72]Tb4b4b4) Tff7b72+ Te6edf3ioDispatcherTb4b4b4)
T8b949e/** Guards against emitting [TAKConnectionEvent.Disconnected] more than once. */
Tff7b72private Tff7b72val Te6edf3disconnectedEmitted Tff7b72= Te6edf3AtomicBooleanTb4b4b4(Tff7b72falseTb4b4b4)
T8b949e/**
* Fail-fast flag checked at the top of [sendXml] so racing broadcasts against a dead connection don't even allocate
* a coroutine.
*/
Tf0883e@Volatile Tff7b72private Tff7b72var Te6edf3closed Tff7b72= Tff7b72false
Tff7b72fun Td2a8ffstartTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3onEventTb4b4b4(Te6edf3TAKConnectionEventTb4b4b4.Te6edf3ConnectedTb4b4b4(Te6edf3currentClientInfoTb4b4b4)Tb4b4b4)
Te6edf3sendProtocolSupportTb4b4b4(Tb4b4b4)
Te6edf3connectionScopeTb4b4b4.Te6edf3launch Tb4b4b4{ Te6edf3readLoopTb4b4b4(Tb4b4b4) Tb4b4b4}
Te6edf3connectionScopeTb4b4b4.Te6edf3launch Tb4b4b4{ Te6edf3keepaliveLoopTb4b4b4(Tb4b4b4) Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffsendProtocolSupportTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3serverUid Tff7b72= Ta5d6ff"Ta5d6ffMeshtastic-TAK-Server-Tffd700${Te6edf3RandomTb4b4b4.Te6edf3nextIntTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3toStringTb4b4b4(Te6edf3TAK_HEX_RADIXTb4b4b4)Tffd700}Ta5d6ff"
Tff7b72val Te6edf3now Tff7b72= Te6edf3ClockTb4b4b4.Te6edf3SystemTb4b4b4.Te6edf3nowTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3stale Tff7b72= Te6edf3now Tff7b72+ Te6edf3TAK_KEEPALIVE_INTERVAL_MSTb4b4b4.Te6edf3milliseconds
Tff7b72val Te6edf3detail Tff7b72=
Ta5d6ff"""
Ta5d6ff <TakControl>
<TakProtocolSupport version=Ta5d6ff"Ta5d6ff0Ta5d6ff"Ta5d6ff/>
</TakControl>
Ta5d6ff"""
Tb4b4b4.Te6edf3trimIndentTb4b4b4(Tb4b4b4)
Te6edf3sendXmlInternalTb4b4b4(Te6edf3buildEventXmlTb4b4b4(Te6edf3uid Tff7b72= Te6edf3serverUidTb4b4b4, Te6edf3type Tff7b72= Ta5d6ff"Ta5d6fft-x-takp-vTa5d6ff"Tb4b4b4, Te6edf3now Tff7b72= Te6edf3nowTb4b4b4, Te6edf3stale Tff7b72= Te6edf3staleTb4b4b4, Te6edf3detail Tff7b72= Te6edf3detailTb4b4b4)Tb4b4b4)
Tb4b4b4}
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffreadLoopTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72try Tb4b4b4{
Tff7b72val Te6edf3buffer Tff7b72= Te6edf3ByteArrayTb4b4b4(Te6edf3TAK_XML_READ_BUFFER_SIZETb4b4b4)
Tff7b72while Tb4b4b4(Te6edf3connectionScopeTb4b4b4.Te6edf3coroutineIsActive Tff7b72&Tff7b72& Tff7b72!Te6edf3closed Tff7b72&Tff7b72& Tff7b72!Te6edf3socketTb4b4b4.Te6edf3isClosedTb4b4b4) Tb4b4b4{
T8b949e// Blocking read off the TLS input stream β must run on the IO dispatcher.
Tff7b72val Te6edf3bytesRead Tff7b72= Te6edf3withContextTb4b4b4(Te6edf3ioDispatcherTb4b4b4) Tb4b4b4{ Te6edf3inputStreamTb4b4b4.Te6edf3readTb4b4b4(Te6edf3bufferTb4b4b4) Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3bytesRead Tff7b72> T79c0ff0Tb4b4b4) Tb4b4b4{
Te6edf3processReceivedDataTb4b4b4(Te6edf3bufferTb4b4b4.Te6edf3copyOfRangeTb4b4b4(T79c0ff0Tb4b4b4, Te6edf3bytesReadTb4b4b4)Tb4b4b4)
Tb4b4b4} Tff7b72else Tff7b72if Tb4b4b4(Te6edf3bytesRead Tff7b72=Tff7b72= Tff7b72-T79c0ff1Tb4b4b4) Tb4b4b4{
Tff7b72break T8b949e// EOF: remote peer closed the connection cleanly
Tb4b4b4}
Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3e
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3closedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3eTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffTAK client read error: Tffd700${Te6edf3currentClientInfoTb4b4b4.Te6edf3idTffd700}Ta5d6ff" Tb4b4b4}
Te6edf3emitDisconnectedTb4b4b4(Te6edf3TAKConnectionEventTb4b4b4.Te6edf3ErrorTb4b4b4(Te6edf3eTb4b4b4)Tb4b4b4)
Tb4b4b4}
Tff7b72return
Tb4b4b4}
Te6edf3emitDisconnectedTb4b4b4(Te6edf3TAKConnectionEventTb4b4b4.Te6edf3DisconnectedTb4b4b4)
Tb4b4b4}
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffkeepaliveLoopTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72while Tb4b4b4(Te6edf3connectionScopeTb4b4b4.Te6edf3coroutineIsActive Tff7b72&Tff7b72& Tff7b72!Te6edf3closed Tff7b72&Tff7b72& Tff7b72!Te6edf3socketTb4b4b4.Te6edf3isClosedTb4b4b4) Tb4b4b4{
Te6edf3kotlinxTb4b4b4.Te6edf3coroutinesTb4b4b4.Te6edf3delayTb4b4b4(Te6edf3TAK_KEEPALIVE_INTERVAL_MSTb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3closedTb4b4b4) Tff7b72break
Te6edf3sendKeepaliveTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffsendKeepaliveTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3now Tff7b72= Te6edf3ClockTb4b4b4.Te6edf3SystemTb4b4b4.Te6edf3nowTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3stale Tff7b72= Te6edf3now Tff7b72+ Te6edf3TAK_KEEPALIVE_INTERVAL_MSTb4b4b4.Te6edf3milliseconds
Te6edf3sendXmlInternalTb4b4b4(Te6edf3buildEventXmlTb4b4b4(Te6edf3uid Tff7b72= Ta5d6ff"Ta5d6fftakPongTa5d6ff"Tb4b4b4, Te6edf3type Tff7b72= Ta5d6ff"Ta5d6fft-x-d-dTa5d6ff"Tb4b4b4, Te6edf3now Tff7b72= Te6edf3nowTb4b4b4, Te6edf3stale Tff7b72= Te6edf3staleTb4b4b4, Te6edf3detail Tff7b72= Ta5d6ff"Ta5d6ff"Tb4b4b4)Tb4b4b4)
Tb4b4b4}
T8b949e/** Respond to ATAK's `t-x-c-t` ping with a pong to reset its RX timeout. */
Tff7b72private Tff7b72fun Td2a8ffsendPongTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3now Tff7b72= Te6edf3ClockTb4b4b4.Te6edf3SystemTb4b4b4.Te6edf3nowTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3stale Tff7b72= Te6edf3now Tff7b72+ Te6edf3TAK_KEEPALIVE_INTERVAL_MSTb4b4b4.Te6edf3milliseconds
Te6edf3sendXmlInternalTb4b4b4(Te6edf3buildEventXmlTb4b4b4(Te6edf3uid Tff7b72= Ta5d6ff"Ta5d6fftakPongTa5d6ff"Tb4b4b4, Te6edf3type Tff7b72= Ta5d6ff"Ta5d6fft-x-c-t-rTa5d6ff"Tb4b4b4, Te6edf3now Tff7b72= Te6edf3nowTb4b4b4, Te6edf3stale Tff7b72= Te6edf3staleTb4b4b4, Te6edf3detail Tff7b72= Ta5d6ff"Ta5d6ff"Tb4b4b4)Tb4b4b4)
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffprocessReceivedDataTb4b4b4(Te6edf3newDataTb4b4b4: Te6edf3ByteArrayTb4b4b4) Tb4b4b4{
Te6edf3frameBufferTb4b4b4.Te6edf3appendTb4b4b4(Te6edf3newDataTb4b4b4)Tb4b4b4.Te6edf3forEach Tb4b4b4{ Te6edf3xmlString Tff7b72-Tff7b72> Te6edf3parseAndHandleMessageTb4b4b4(Te6edf3xmlStringTb4b4b4) Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffparseAndHandleMessageTb4b4b4(Te6edf3xmlStringTb4b4b4: Tffa657StringTb4b4b4) Tb4b4b4{
T8b949e// Fast-path: detect keepalive pings before full XML parsing to avoid
T8b949e// both the parse overhead and the noisy RAW CoT IN log line every 4.5s.
Tff7b72if Tb4b4b4(Te6edf3xmlStringTb4b4b4.Te6edf3containsTb4b4b4(Ta5d6ff"Ta5d6fft-x-c-tTa5d6ff"Tb4b4b4) Tff7b72|Tff7b72| Te6edf3xmlStringTb4b4b4.Te6edf3containsTb4b4b4(Ta5d6ff"Ta5d6ffuid=Ta5d6ff\"Ta5d6ffpingTa5d6ff\"Ta5d6ff"Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3sendPongTb4b4b4(Tb4b4b4)
Tff7b72return
Tb4b4b4}
T8b949e// Full raw CoT XML from the ATAK client, before any parsing happens.
T8b949e// Emitted at debug level so it's always available in logcat for field
T8b949e// debugging without needing a release rebuild. Not truncated β the
T8b949e// reader of this log needs the complete event to reproduce issues.
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffRAW CoT IN (TCP Tffd700${Te6edf3currentClientInfoTb4b4b4.Te6edf3idTffd700}Ta5d6ff): Tffd700$Te6edf3xmlStringTa5d6ff" Tb4b4b4}
Tff7b72val Te6edf3parser Tff7b72= Te6edf3CoTXmlParserTb4b4b4(Te6edf3xmlStringTb4b4b4)
Tff7b72val Te6edf3result Tff7b72= Te6edf3parserTb4b4b4.Te6edf3parseTb4b4b4(Tb4b4b4)
Te6edf3result
Tb4b4b4.Te6edf3onSuccess Tb4b4b4{ Te6edf3cotMessage Tff7b72-Tff7b72>
Tff7b72when Tb4b4b4{
Te6edf3cotMessageTb4b4b4.Te6edf3typeTb4b4b4.Te6edf3startsWithTb4b4b4(Ta5d6ff"Ta5d6fft-x-takpTa5d6ff"Tb4b4b4) Tff7b72-Tff7b72> Tb4b4b4{
Te6edf3handleProtocolControlTb4b4b4(Te6edf3cotMessageTb4b4b4.Te6edf3typeTb4b4b4, Te6edf3xmlStringTb4b4b4)
Tff7b72return
Tb4b4b4}
Tff7b72else Tff7b72-Tff7b72> Tb4b4b4{
Te6edf3cotMessageTb4b4b4.Te6edf3contactTff7b72?.Te6edf3let Tb4b4b4{ Te6edf3contact Tff7b72-Tff7b72>
Tff7b72val Te6edf3updatedClientInfo Tff7b72=
Te6edf3currentClientInfoTb4b4b4.Te6edf3copyTb4b4b4(
Te6edf3callsign Tff7b72= Te6edf3currentClientInfoTb4b4b4.Te6edf3callsign Tff7b72?: Te6edf3contactTb4b4b4.Te6edf3callsignTb4b4b4,
Te6edf3uid Tff7b72= Te6edf3currentClientInfoTb4b4b4.Te6edf3uid Tff7b72?: Te6edf3cotMessageTb4b4b4.Te6edf3uidTb4b4b4,
Tb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3updatedClientInfo Tff7b72!Tff7b72= Te6edf3currentClientInfoTb4b4b4) Tb4b4b4{
Te6edf3currentClientInfo Tff7b72= Te6edf3updatedClientInfo
Te6edf3onEventTb4b4b4(Te6edf3TAKConnectionEventTb4b4b4.Te6edf3ClientInfoUpdatedTb4b4b4(Te6edf3updatedClientInfoTb4b4b4)Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Te6edf3onEventTb4b4b4(Te6edf3TAKConnectionEventTb4b4b4.Te6edf3MessageTb4b4b4(Te6edf3cotMessageTb4b4b4, Te6edf3currentClientInfoTb4b4b4)Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4.Te6edf3onFailure Tb4b4b4{ Te6edf3e Tff7b72-Tff7b72> Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3eTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to parse CoT XML from TAK client Tffd700${Te6edf3currentClientInfoTb4b4b4.Te6edf3idTffd700}Ta5d6ff" Tb4b4b4} Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffhandleProtocolControlTb4b4b4(Te6edf3typeTb4b4b4: Tffa657StringTb4b4b4, Te6edf3xmlStringTb4b4b4: Tffa657StringTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3type Tff7b72=Tff7b72= Ta5d6ff"Ta5d6fft-x-takp-qTa5d6ff"Tb4b4b4) Tb4b4b4{
Te6edf3sendProtocolResponseTb4b4b4(Tb4b4b4)
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffUnhandled protocol control type: Tffd700$Te6edf3typeTa5d6ff (raw=Tffd700$Te6edf3xmlStringTa5d6ff)Ta5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffsendProtocolResponseTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3serverUid Tff7b72= Ta5d6ff"Ta5d6ffMeshtastic-TAK-Server-Tffd700${Te6edf3RandomTb4b4b4.Te6edf3nextIntTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3toStringTb4b4b4(Te6edf3TAK_HEX_RADIXTb4b4b4)Tffd700}Ta5d6ff"
Tff7b72val Te6edf3now Tff7b72= Te6edf3ClockTb4b4b4.Te6edf3SystemTb4b4b4.Te6edf3nowTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3stale Tff7b72= Te6edf3now Tff7b72+ Te6edf3TAK_KEEPALIVE_INTERVAL_MSTb4b4b4.Te6edf3milliseconds
Tff7b72val Te6edf3detail Tff7b72=
Ta5d6ff"""
Ta5d6ff <TakControl>
<TakResponse status=Ta5d6ff"Ta5d6fftrueTa5d6ff"Ta5d6ff/>
</TakControl>
Ta5d6ff"""
Tb4b4b4.Te6edf3trimIndentTb4b4b4(Tb4b4b4)
Te6edf3sendXmlInternalTb4b4b4(Te6edf3buildEventXmlTb4b4b4(Te6edf3uid Tff7b72= Te6edf3serverUidTb4b4b4, Te6edf3type Tff7b72= Ta5d6ff"Ta5d6fft-x-takp-rTa5d6ff"Tb4b4b4, Te6edf3now Tff7b72= Te6edf3nowTb4b4b4, Te6edf3stale Tff7b72= Te6edf3staleTb4b4b4, Te6edf3detail Tff7b72= Te6edf3detailTb4b4b4)Tb4b4b4)
Tb4b4b4}
Tff7b72fun Td2a8ffsendTb4b4b4(Te6edf3cotMessageTb4b4b4: Te6edf3CoTMessageTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3closedTb4b4b4) Tff7b72return
Tff7b72val Te6edf3xml Tff7b72= Te6edf3cotMessageTb4b4b4.Te6edf3toXmlTb4b4b4(Tb4b4b4)
T8b949e// Full raw CoT XML being shipped out to the ATAK client, after the
T8b949e// CoTMessage β XML round trip. This is the exact bytes the client
T8b949e// will receive, so logging here closes the debugging loop with the
T8b949e// matching RAW CoT IN line on the receiver.
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffRAW CoT OUT (TCP Tffd700${Te6edf3currentClientInfoTb4b4b4.Te6edf3idTffd700}Ta5d6ff): Tffd700$Te6edf3xmlTa5d6ff" Tb4b4b4}
Te6edf3sendXmlInternalTb4b4b4(Te6edf3xmlTb4b4b4)
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffbuildEventXmlTb4b4b4(Te6edf3uidTb4b4b4: Tffa657StringTb4b4b4, Te6edf3typeTb4b4b4: Tffa657StringTb4b4b4, Te6edf3nowTb4b4b4: Te6edf3InstantTb4b4b4, Te6edf3staleTb4b4b4: Te6edf3InstantTb4b4b4, Te6edf3detailTb4b4b4: Tffa657StringTb4b4b4)Tb4b4b4: Tffa657String Tb4b4b4{
Tff7b72val Te6edf3detailContent Tff7b72= Tff7b72if Tb4b4b4(Te6edf3detailTb4b4b4.Te6edf3isBlankTb4b4b4(Tb4b4b4)Tb4b4b4) Ta5d6ff"Ta5d6ff<detail/>Ta5d6ff" Tff7b72else Ta5d6ff"Ta5d6ff<detail>Tffd700$Te6edf3detailTa5d6ff</detail>Ta5d6ff"
Tff7b72val Te6edf3point Tff7b72= Ta5d6ff"""Ta5d6ff<point lat=Ta5d6ff"Ta5d6ff0Ta5d6ff"Ta5d6ff lon=Ta5d6ff"Ta5d6ff0Ta5d6ff"Ta5d6ff hae=Ta5d6ff"Ta5d6ff0Ta5d6ff"Ta5d6ff ce=Ta5d6ff"Tffd700$Te6edf3TAK_UNKNOWN_POINT_VALUETa5d6ff"Ta5d6ff le=Ta5d6ff"Tffd700$Te6edf3TAK_UNKNOWN_POINT_VALUETa5d6ff"Ta5d6ff/>Ta5d6ff"""
Tff7b72return Ta5d6ff"""Ta5d6ff<event version=Ta5d6ff"Ta5d6ff2.0Ta5d6ff"Ta5d6ff uid=Ta5d6ff"Tffd700$Te6edf3uidTa5d6ff"Ta5d6ff type=Ta5d6ff"Tffd700$Te6edf3typeTa5d6ff"Ta5d6ff time=Ta5d6ff"Tffd700$Te6edf3nowTa5d6ff"Ta5d6ff start=Ta5d6ff"Tffd700$Te6edf3nowTa5d6ff"Ta5d6ff stale=Ta5d6ff"Tffd700$Te6edf3staleTa5d6ff"Ta5d6ff how=Ta5d6ff"Ta5d6ffm-gTa5d6ff"Ta5d6ff>Ta5d6ff""" Tff7b72+
Te6edf3point Tff7b72+
Te6edf3detailContent Tff7b72+
Ta5d6ff"Ta5d6ff</event>Ta5d6ff"
Tb4b4b4}
T8b949e/**
* Send raw XML directly to this client. Used for mesh-originated messages that bypass CoTMessage parsing to
* preserve shape detail elements.
*/
Tff7b72fun Td2a8ffsendRawXmlTb4b4b4(Te6edf3xmlTb4b4b4: Tffa657StringTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffRAW CoT OUT (TCP Tffd700${Te6edf3currentClientInfoTb4b4b4.Te6edf3idTffd700}Ta5d6ff): [raw] Tffd700$Te6edf3xmlTa5d6ff" Tb4b4b4}
Te6edf3sendXmlInternalTb4b4b4(Te6edf3xmlTb4b4b4)
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffsendXmlInternalTb4b4b4(Te6edf3xmlTb4b4b4: Tffa657StringTb4b4b4) Tb4b4b4{
T8b949e// Fail-fast synchronous check BEFORE allocating a coroutine. This is the hot path
T8b949e// for broadcasts β see the scope doc above for why it matters.
Tff7b72if Tb4b4b4(Te6edf3closedTb4b4b4) Tff7b72return
Te6edf3connectionScopeTb4b4b4.Te6edf3launch Tb4b4b4{
T8b949e// Re-check inside the coroutine: we may have been cancelled or marked closed
T8b949e// between the launch and the dispatcher picking this up.
Tff7b72if Tb4b4b4(Te6edf3closedTb4b4b4) Tff7b72returnTf0883e@launch
Tff7b72try Tb4b4b4{
Te6edf3writeMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3closed Tff7b72|Tff7b72| Te6edf3socketTb4b4b4.Te6edf3isClosedTb4b4b4) Tff7b72returnTf0883e@withLock
Tff7b72val Te6edf3bytes Tff7b72= Te6edf3xmlTb4b4b4.Te6edf3toByteArrayTb4b4b4(Te6edf3CharsetsTb4b4b4.Te6edf3UTF_8Tb4b4b4)
T8b949e// Blocking write on TLS output must run on the IO dispatcher
Te6edf3withContextTb4b4b4(Te6edf3ioDispatcherTb4b4b4) Tb4b4b4{
Te6edf3outputStreamTb4b4b4.Te6edf3writeTb4b4b4(Te6edf3bytesTb4b4b4)
Te6edf3outputStreamTb4b4b4.Te6edf3flushTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3e
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
T8b949e// Don't spam on writes that raced a disconnect we already observed.
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3closedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3eTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffTAK client send error: Tffd700${Te6edf3currentClientInfoTb4b4b4.Te6edf3idTffd700}Ta5d6ff" Tb4b4b4}
Te6edf3emitDisconnectedTb4b4b4(Te6edf3TAKConnectionEventTb4b4b4.Te6edf3ErrorTb4b4b4(Te6edf3eTb4b4b4)Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72fun Td2a8ffcloseTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3frameBufferTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3emitDisconnectedTb4b4b4(Te6edf3TAKConnectionEventTb4b4b4.Te6edf3DisconnectedTb4b4b4)
Tb4b4b4}
T8b949e/**
* Emits [event] (expected to be [TAKConnectionEvent.Disconnected] or [TAKConnectionEvent.Error]) at most once
* across all code paths, then tears down the per-connection coroutines and socket.
*
* This is the ONLY place the connection's entire coroutine scope β keepalive loop, read loop, and any in-flight
* send coroutines β gets cancelled when the *remote* peer closes the TLS stream. Without this, Java's
* [Socket.isClosed] only reports whether *our* side called close(), so the keepalive loop's `!socket.isClosed`
* guard never fires, the broadcast fanout keeps launching writes onto the dead socket via [sendXml], and every
* iteration logs `SSLOutputStream / Socket closed`. Before [closed] + [connectionScope.cancel] were added, a single
* session with a few reconnects accumulated hundreds of zombie write coroutines each spamming errors in parallel.
*
* Idempotent via [AtomicBoolean.compareAndSet], so racing calls from [readLoop], [keepaliveLoop], and [sendXml] all
* converge on a single teardown.
*/
Tff7b72private Tff7b72fun Td2a8ffemitDisconnectedTb4b4b4(Te6edf3eventTb4b4b4: Te6edf3TAKConnectionEventTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3disconnectedEmittedTb4b4b4.Te6edf3compareAndSetTb4b4b4(Tff7b72falseTb4b4b4, Tff7b72trueTb4b4b4)Tb4b4b4) Tb4b4b4{
T8b949e// Set the fail-fast flag BEFORE emitting the event. [TAKServerJvm] will
T8b949e// schedule an async map removal on receipt, and any broadcast racing the
T8b949e// removal must see `closed = true` when it hits [send] / [sendXml].
Te6edf3closed Tff7b72= Tff7b72true
Te6edf3onEventTb4b4b4(Te6edf3eventTb4b4b4)
T8b949e// Cancel the whole scope β readLoop, keepaliveLoop, and every queued or
T8b949e// in-flight sendXml coroutine. Any write blocked in the syscall will throw
T8b949e// on the next iteration because we close the socket next.
Te6edf3connectionScopeTb4b4b4.Te6edf3cancelTb4b4b4(Tb4b4b4)
Tff7b72try Tb4b4b4{
Te6edf3socketTb4b4b4.Te6edf3closeTb4b4b4(Tb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3_Tb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Served by rngit 1.5.0 - Generated in 0.08s